Repository navigation
feat(streaming): add EventStreamer for iterating stream events - #1567
Conversation
|
@jakelorocco here's my initial draft POC for a events iterator for streaming. This is the better of the two POCs I wrote and I'll leave this here for your review while I'm out. Once I get back next week we can sync and I'll continue work on this |
jakelorocco
left a comment
There was a problem hiding this comment.
This pretty much looks like exactly what I was looking for; thank you!
7d03c9a to
52987f8
Compare
stream(as_events=True) returns an EventStreamer that yields a stream's typed StreamEvent objects via async-for, instead of the Streamer's str chunks. It runs the stream to completion on a background task and delivers events through a queue, so the terminal CompletedEvent and ErrorEvent reach the consumer through iteration, staying queued behind an early break until it resumes — which iterating the driver's generator cannot guarantee. Additive and opt-in: the default stream() chunk contract is unchanged. Events are threaded to an optional per-stream queue in _emit_event; the STREAMING_EVENT hook still fires on all paths, so telemetry is unaffected. stream() is split into _stream() plus a dispatching wrapper with typed overloads on as_events. Assisted-by: Claude Code Signed-off-by: Alex Bozarth <ajbozart@us.ibm.com>
52987f8 to
81f392a
Compare
|
PR is updated and moved out of draft POC, ready for review |
planetf1
left a comment
There was a problem hiding this comment.
I found three correctness and documentation issues inline.
… setup Split setup handling on asyncio.CancelledError instead of Exception, so KeyboardInterrupt/SystemExit/custom BaseExceptions surface to the caller rather than a masking CancelledError. Also import the event types the EventStreamer how-to example uses. Assisted-by: Claude Code Signed-off-by: Alex Bozarth <ajbozart@us.ibm.com>
planetf1
left a comment
There was a problem hiding this comment.
All review threads resolved and verified against 6a65309:
- BaseException preservation during setup: verified — pump splits on CancelledError vs other exceptions, original error surfaces to the caller, regression test updated to expect it.
- Missing event-type imports in the how-to example: verified.
- Bounded-queue proposal: withdrawn — unbounded eager buffering is the right trade (measured ~160 B/event, O(output length) per stream; a blocking put would pin the backend slot and breaks the close path).
CI green.
|
@jakelorocco if you could give this another review pass before I merge, I want to make sure the edits I made to the POC after your initial approval are still good |
jakelorocco
left a comment
There was a problem hiding this comment.
Found one potential issue; I don't know if it's worth addressing / worth holding up this PR; will defer to you
| if self._pump_task is not None: | ||
| if not self._pump_task.done(): | ||
| self._pump_task.cancel() | ||
| try: | ||
| await self._pump_task | ||
| except asyncio.CancelledError: | ||
| pass # we cancelled it; consume the resulting CancelledError |
There was a problem hiding this comment.
I think the assumption behind this exception block may be incorrect. It could also be raised if the current task (the one running this aclose) is cancelled from outside of it. In that case, I think we ought to be re-raising the cancelled err.
I don't know how important this is, but I think mots already handle it with their .cancel_generation field. Maybe we can do something similar here?
There was a problem hiding this comment.
updated to match how MOT handles it and added a test like MOTs as well in cba531a
…lose aclose swallowed every CancelledError while awaiting the cancelled pump, absorbing an external cancel of the task running aclose. Re-raise when current_task().cancelling() > 0, mirroring MOT.cancel_generation. Assisted-by: Claude Code Signed-off-by: Alex Bozarth <ajbozart@us.ibm.com>
0c02eae
Issue
Follow-up to the design discussion on #1543 (now merged) — streaming event consumption ergonomics.
Description
Adds an opt-in
EventStreamer, returned bystream(..., as_events=True), for iterating a single stream's typedStreamEventobjects directly withasync for. The only other route to those events is the process-globalSTREAMING_EVENThook, which is awkward for a single stream — you register a plugin and demultiplex onstreaming_id.Under the hood it wraps a
Streamer, drives it to completion on a background task, and delivers the events through a queue the consumer drains withasync for.It is additive and opt-in:
stream()'s default chunk-iteratingStreameris unchanged, and theSTREAMING_EVENThook still fires on every stream, so telemetry and multi-stream observers are unaffected.EventStreameris just the ergonomic way to consume one stream's events inline; the hook remains the right tool for observing many at once.The streaming examples and docs go back to iterating events, reverting the
streaming_event-hook workaround #1543 adopted when the iterator was removed; the v0.7→v0.8 migration guide now mapsresult.events()tostream(as_events=True)to match.Testing
Unit tests cover
EventStreamer's event-mode paths (natural completion, requirement failure, error re-raise, early break), resume-and-drain after a break, teardown edge cases (idempotentaclose, a faulted/abandoned pump, aBaseExceptionduring setup that must not hang), and setup-error surfacing;test/typing/check_stream_as_events.pycovers theas_eventsoverloads including a dynamicbool.Attribution
Adding a new component, requirement, sampling strategy, or tool?
If your PR adds or modifies one of the types below, check the matching box. A checklist of type-specific review items will be posted as a comment.
NOTE: Please ensure you have an issue that has been acknowledged by a core contributor and routed you to open a pull request against this repository. Otherwise, please open an issue before continuing with this pull request.